🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / network / src / commonMain / kotlin / org / meshtastic / core / network / radio / StreamTransport.kt
Displaying Raw • Download
core/network/src/commonMain/kotlin/org/meshtastic/core/network/radio/StreamTransport.kt bd2863243bab6eb213401d949839a2bc74dde7e2 (bd286324) Text, 8.18 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.network.radio
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.cancelAndJoin
Tff7b72import T7ee787kotlinx.coroutines.channels.Channel
Tff7b72import T7ee787kotlinx.coroutines.currentCoroutineContext
Tff7b72import T7ee787kotlinx.coroutines.isActive
Tff7b72import T7ee787kotlinx.coroutines.withTimeoutOrNull
Tff7b72import T7ee787org.meshtastic.core.common.util.handledLaunch
Tff7b72import T7ee787org.meshtastic.core.network.transport.StreamFrameCodec
Tff7b72import T7ee787org.meshtastic.core.repository.RadioTransport
Tff7b72import T7ee787org.meshtastic.core.repository.RadioTransportCallback
Tff7b72import T7ee787org.meshtastic.core.repository.TransportDisconnectReason
Tff7b72import T7ee787kotlin.time.Duration.Companion.seconds
Tff7b72private Tff7b72val Te6edf3STREAM_CLOSE_DRAIN_TIMEOUT Tff7b72= T79c0ff1T79c0ff0.Te6edf3seconds
T8b949e/**
* An interface that assumes we are talking to a meshtastic device over some sort of stream connection (serial or TCP
* probably).
*
* Delegates framing logic to [StreamFrameCodec] from `core:network`.
*/
Tff7b72abstract Tff7b72class T56d364StreamTransportTb4b4b4(Tff7b72protected Tff7b72val Te6edf3callbackTb4b4b4: Te6edf3RadioTransportCallbackTb4b4b4, Tff7b72protected Tff7b72val Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4) Tb4b4b4:
Te6edf3RadioTransport Tb4b4b4{
Tff7b72private Tff7b72class T56d364FramedSendTb4b4b4(
Tff7b72val Te6edf3payloadTb4b4b4: Te6edf3ByteArrayTb4b4b4,
Tff7b72val Te6edf3writerTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3ByteArrayTb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4,
Tff7b72val Te6edf3flusherTb4b4b4: Te6edf3suspend Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4,
Tff7b72val Te6edf3onCompletionTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4,
Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3codec Tff7b72=
Te6edf3StreamFrameCodecTb4b4b4(Te6edf3onPacketReceived Tff7b72= Tb4b4b4{ Te6edf3callbackTb4b4b4.Te6edf3handleFromRadioTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}Tb4b4b4, Te6edf3logTag Tff7b72= Ta5d6ff"Ta5d6ffStreamTransportTa5d6ff"Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3sendQueue Tff7b72= Te6edf3ChannelTff7b72<Te6edf3FramedSendTff7b72>Tb4b4b4(Te6edf3capacity Tff7b72= Te6edf3MAX_PENDING_SENDSTb4b4b4)
Tff7b72private Tff7b72val Te6edf3sendWorker Tff7b72=
Te6edf3scope
Tb4b4b4.Te6edf3handledLaunch Tb4b4b4{
Tff7b72for Tb4b4b4(Te6edf3send Tff7b72in Te6edf3sendQueueTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3failure Tff7b72=
Te6edf3runCatching Tb4b4b4{ Te6edf3codecTb4b4b4.Te6edf3frameAndSendTb4b4b4(Te6edf3sendTb4b4b4.Te6edf3payloadTb4b4b4, Te6edf3sendTb4b4b4.Te6edf3writerTb4b4b4, Te6edf3sendTb4b4b4.Te6edf3flusherTb4b4b4) Tb4b4b4}Tb4b4b4.Te6edf3exceptionOrNullTb4b4b4(Tb4b4b4)
Tff7b72try Tb4b4b4{
Tff7b72when Tb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{
Tff7b72null Tff7b72-Tff7b72> Tffa657Unit
Tff7b72is Te6edf3CancellationException Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3currentCoroutineContextTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3isActiveTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3failure
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffStreamTransport: framed send cancelledTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72is Te6edf3Error Tff7b72-Tff7b72> Tff7b72throw Te6edf3failure
Tff7b72else Tff7b72-Tff7b72> Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffStreamTransport: framed send failedTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4} Tff7b72finally Tb4b4b4{
Te6edf3completeSendTb4b4b4(Te6edf3sendTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4.Te6edf3also Tb4b4b4{ Te6edf3worker Tff7b72-Tff7b72> Te6edf3workerTb4b4b4.Te6edf3invokeOnCompletion Tb4b4b4{ Te6edf3drainAbandonedSendsTb4b4b4(Tb4b4b4) Tb4b4b4} Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffdrainAbandonedSendsTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3sendQueueTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4)
Tff7b72while Tb4b4b4(Tff7b72trueTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3abandoned Tff7b72= Te6edf3sendQueueTb4b4b4.Te6edf3tryReceiveTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3getOrNullTb4b4b4(Tb4b4b4) Tff7b72?: Tff7b72break
Te6edf3completeSendTb4b4b4(Te6edf3abandonedTb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Closes admission and lets already-admitted frames finish before cancelling the worker as a bounded fallback.
*
* A forced cancellation can interrupt a frame mid-write, so [close] is terminal for the underlying stream: a
* subclass must close or replace that stream before another transport instance writes to it.
*/
Tff7b72override Tff7b72suspend Tff7b72fun Td2a8ffcloseTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3d Tb4b4b4{ Ta5d6ff"Ta5d6ffClosing stream transportTa5d6ff" Tb4b4b4}
Te6edf3sendQueueTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4)
Tff7b72val Te6edf3drained Tff7b72= Te6edf3withTimeoutOrNullTb4b4b4(Te6edf3STREAM_CLOSE_DRAIN_TIMEOUTTb4b4b4) Tb4b4b4{ Te6edf3sendWorkerTb4b4b4.Te6edf3joinTb4b4b4(Tb4b4b4) Tb4b4b4} Tff7b72!Tff7b72= Tff7b72null
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3drainedTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffStreamTransport: framed send drain timed out after Tffd700$Te6edf3STREAM_CLOSE_DRAIN_TIMEOUTTa5d6ff" Tb4b4b4}
Te6edf3sendWorkerTb4b4b4.Te6edf3cancelAndJoinTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Signals the transport callback that the device has disconnected and optionally waits for the transport to stop.
*
* @param waitForStopped if true we should wait for the transport to finish - must be false if called from inside
* transport callbacks
* @param isPermanent true when the user explicitly disconnects (e.g. [close] was called), or when an authorization
* failure makes the current connection attempt unrecoverable. USB unplug, I/O errors, and similar conditions are
* transient — the transport may recover when the device is replugged or the OS re-enumerates. Defaults to false
* so callbacks default to "may come back".
* @param errorMessage optional user-facing reason for the disconnect; surfaced via
* [RadioTransportCallback.onDisconnect]. Prefer [reason] for newly-classified disconnect causes so transports do
* not own UI copy.
* @param reason optional structured cause for service/UI-specific handling.
*/
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffonDeviceDisconnectTb4b4b4(
Te6edf3waitForStoppedTb4b4b4: Tffa657BooleanTb4b4b4,
Te6edf3isPermanentTb4b4b4: Tffa657Boolean Tff7b72= Tff7b72falseTb4b4b4,
Te6edf3errorMessageTb4b4b4: Tffa657String? Tff7b72= Tff7b72nullTb4b4b4,
Te6edf3reasonTb4b4b4: Te6edf3TransportDisconnectReason? Tff7b72= Tff7b72nullTb4b4b4,
Tb4b4b4) Tb4b4b4{
Te6edf3callbackTb4b4b4.Te6edf3onDisconnectTb4b4b4(Te6edf3isPermanent Tff7b72= Te6edf3isPermanentTb4b4b4, Te6edf3errorMessage Tff7b72= Te6edf3errorMessageTb4b4b4, Te6edf3reason Tff7b72= Te6edf3reasonTb4b4b4)
Tb4b4b4}
Tff7b72protected Tff7b72open Tff7b72fun Td2a8ffconnectTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Before connecting, send a few START1s to wake a sleeping device
Te6edf3sendBytesTb4b4b4(Te6edf3StreamFrameCodecTb4b4b4.Te6edf3WAKE_BYTESTb4b4b4)
T8b949e// Now tell clients they can (finally use the api)
Te6edf3callbackTb4b4b4.Te6edf3onConnectTb4b4b4(Tb4b4b4)
Tb4b4b4}
T8b949e/**
* Writes raw bytes to the underlying stream (serial port, TCP socket, etc.). Implementations may block until the
* driver accepts the bytes, so direct callers must already be running on an I/O dispatcher. Framed packet sends use
* [queueFramedSend], whose writer controls dispatcher confinement. Subclasses that rely on [queueFramedSend]'s
* default writer must therefore supply an I/O-confined [scope].
*/
Tff7b72abstract Tff7b72fun Td2a8ffsendBytesTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ByteArrayTb4b4b4)
T8b949e/** Flushes buffered bytes to the underlying stream. No-op by default. */
Tff7b72open Tff7b72fun Td2a8ffflushBytesTb4b4b4(Tb4b4b4) Tb4b4b4{Tb4b4b4}
T8b949e/**
* Queues the framed packet onto [scope], optionally binding the deferred write to a transport-session resource.
*
* Capturing the writer at admission keeps a queued write from being redirected to a replacement connection before
* its coroutine runs. The result reports only successful scheduling on a live scope; physical delivery is not
* confirmed.
*
* [onCompletion] runs exactly once whether the send is admitted, abandoned during worker teardown, or rejected.
*/
Tff7b72protected Tff7b72fun Td2a8ffqueueFramedSendTb4b4b4(
Te6edf3payloadTb4b4b4: Te6edf3ByteArrayTb4b4b4,
Te6edf3writerTb4b4b4: Te6edf3suspend Tb4b4b4(Te6edf3ByteArrayTb4b4b4) Tff7b72-Tff7b72> Tffa657Unit Tff7b72= Tb4b4b4{ Te6edf3sendBytesTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}Tb4b4b4,
Te6edf3flusherTb4b4b4: Te6edf3suspend Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657Unit Tff7b72= Tb4b4b4{ Te6edf3flushBytesTb4b4b4(Tb4b4b4) Tb4b4b4}Tb4b4b4,
Te6edf3onCompletionTb4b4b4: Tb4b4b4(Tb4b4b4) Tff7b72-Tff7b72> Tffa657Unit Tff7b72= Tb4b4b4{Tb4b4b4}Tb4b4b4,
Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3send Tff7b72= Te6edf3FramedSendTb4b4b4(Te6edf3payloadTb4b4b4, Te6edf3writerTb4b4b4, Te6edf3flusherTb4b4b4, Te6edf3onCompletionTb4b4b4)
Tff7b72val Te6edf3accepted Tff7b72= Te6edf3sendWorkerTb4b4b4.Te6edf3isActive Tff7b72&Tff7b72& Te6edf3sendQueueTb4b4b4.Te6edf3trySendTb4b4b4(Te6edf3sendTb4b4b4)Tb4b4b4.Te6edf3isSuccess
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3acceptedTb4b4b4) Te6edf3completeSendTb4b4b4(Te6edf3sendTb4b4b4)
Tff7b72return Te6edf3accepted
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffcompleteSendTb4b4b4(Te6edf3sendTb4b4b4: Te6edf3FramedSendTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3failure Tff7b72= Te6edf3runCatchingTb4b4b4(Te6edf3sendTb4b4b4.Te6edf3onCompletionTb4b4b4)Tb4b4b4.Te6edf3exceptionOrNullTb4b4b4(Tb4b4b4)
Tff7b72when Tb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{
Tff7b72null Tff7b72-Tff7b72> Tffa657Unit
Tff7b72is Te6edf3Error Tff7b72-Tff7b72> Tff7b72throw Te6edf3failure
Tff7b72else Tff7b72-Tff7b72> Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3failureTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffStreamTransport: framed-send completion failedTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffhandleSendToRadioTb4b4b4(Te6edf3pTb4b4b4: Te6edf3ByteArrayTb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3queueFramedSendTb4b4b4(Te6edf3pTb4b4b4)
T8b949e/** Process a single incoming byte through the stream framing state machine. */
Tff7b72protected Tff7b72fun Td2a8ffreadCharTb4b4b4(Te6edf3cTb4b4b4: Tffa657ByteTb4b4b4) Tb4b4b4{
Te6edf3codecTb4b4b4.Te6edf3processInputByteTb4b4b4(Te6edf3cTb4b4b4)
Tb4b4b4}
Tff7b72internal Tff7b72companion Tff7b72object Tb4b4b4{
T8b949e/** Upper bound on framed sends awaiting the serialized writer. */
Tff7b72const Tff7b72val Te6edf3MAX_PENDING_SENDS Tff7b72= T79c0ff3T79c0ff2
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.0 - Generated in 0.06s